Java NIO - Netty 支持 WebSocket 详细案例以及深度优化方案

从简单案例开始

WebSocket 协议在握手阶段使用的是 HTTP 协议,握手成功后会升级(Upgrade)为双向通道,数据传输不再使用普通的字节流,而是使用名为 Frame(数据帧)的基本单位。Netty 贴心地为我们提供了很多开箱即用的编解码器,使得我们只需要几行配置就能搭建好它。以下是完整的服务端代码实现:


主引导类的编写

WebSocket 的核心在于它的 Pipeline(责任链)流水线配置。我们需要添加 HTTP 编解码器、聚合器,以及 Netty 自带的 WebSocket 协议服务器协议处理器。

MyWebSocketServer:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelInitializer;
import io.netty.channel.ChannelPipeline;
import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.SocketChannel;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import io.netty.handler.codec.http.HttpObjectAggregator;
import io.netty.handler.codec.http.HttpServerCodec;
import io.netty.handler.codec.http.websocketx.WebSocketServerProtocolHandler;
import io.netty.handler.stream.ChunkedWriteHandler;

public class MyWebSocketServer {

public static void main(String[] args) {
EventLoopGroup bossGroup = new NioEventLoopGroup(1);
EventLoopGroup workerGroup = new NioEventLoopGroup();
try {
ServerBootstrap serverBootstrap = new ServerBootstrap();
serverBootstrap.group(bossGroup, workerGroup)
.channel(NioServerSocketChannel.class)
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ChannelPipeline pipeline = ch.pipeline();

// 1. 因为基于 Http 协议进行握手,所以必须添加 HttpServerCodec(Http编解码器)
pipeline.addLast(new HttpServerCodec());

// 插入跨域配置
/*CorsConfig corsConfig = CorsConfigBuilder.forAnyOrigin()
.allowNullOrigin() // 允许本地 file:// 协议带来的 null 跨域
.allowCredentials()
.build();
pipeline.addLast(new CorsHandler(corsConfig));*/

// 2. 以块传输的方式写数据,主要用于支持传输大文件或流数据
pipeline.addLast(new ChunkedWriteHandler());

// 3. Http 数据在传输过程中是分段的,HttpObjectAggregator 能够将多个段聚合起来
// 参数 8192 代表聚合后的最大内容长度(8KB)
pipeline.addLast(new HttpObjectAggregator(8192));

/*
* 4. 核心:WebSocket 协议服务器处理器
* - 它会指定客户端连接的 URL 路由(比如下面的 "/hello")
* - 它会自动处理 WebSocket 握手状态切换 (将 HTTP 升级为 WebSocket)
* - 它会自动处理心跳控制帧(Ping / Pong)
*/
pipeline.addLast(new WebSocketServerProtocolHandler("/hello"));

// 5. 挂载我们自定义的文本帧处理器
pipeline.addLast(new MyWebSocketHandler());
}
});
System.out.println("WebSocket 服务端正在启动...");
ChannelFuture channelFuture = serverBootstrap.bind(8888).sync();
System.out.println("WebSocket 服务端启动成功,监听端口:" + 8888);
channelFuture.channel().closeFuture().sync();
} catch (InterruptedException e) {
throw new RuntimeException(e);
} finally {
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully();
}
}
}


业务处理器的编写

MyWebSocketHandler:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.SimpleChannelInboundHandler;
import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;
import java.time.LocalDateTime;

public class MyWebSocketHandler extends SimpleChannelInboundHandler<TextWebSocketFrame> {

@Override
protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception {
// 1. 获取客户端传过来的文本内容
String text = msg.text();
System.out.println("【WebSocket服务端】收到来自客户端[" + ctx.channel().remoteAddress() + "]的消息: " + text);

// 2. 回复客户端一条消息(注意:回复也必须包装成 TextWebSocketFrame 对象)
String responseText = "【服务器回复】" + LocalDateTime.now() + " -> 我已收到你的消息: " + text;
ctx.channel().writeAndFlush(new TextWebSocketFrame(responseText));
}

@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
System.out.println("【WebSocket服务端】客户端连接成功: " + ctx.channel().remoteAddress());
}

@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
System.out.println("【WebSocket服务端】客户端断开连接: " + ctx.channel().remoteAddress());
}

@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
System.err.println("【WebSocket服务端】发生异常" + cause);
ctx.close(); // 关闭通道
}
}


浏览器客户端测试

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
<!DOCTYPE html>
<html lang="zh-CN">
<head>
<meta charset="UTF-8">
<title>Netty WebSocket 测试客户端</title>
<style>
body {
font-family: 'Segoe UI', Tahoma, Geneva, Verdana, sans-serif;
background-color: #f4f7f6;
margin: 0;
padding: 20px;
display: flex;
justify-content: center;
}
.container {
width: 100%;
max-width: 600px;
background: #ffffff;
padding: 25px;
border-radius: 10px;
box-shadow: 0 4px 15px rgba(0,0,0,0.08);
}
h2 {
margin-top: 0;
color: #333;
border-bottom: 2px solid #eaeaea;
padding-bottom: 10px;
}
.status-container {
display: flex;
align-items: center;
margin-bottom: 15px;
font-size: 14px;
}
.status-dot {
width: 10px;
height: 10px;
border-radius: 50%;
background-color: #ff4d4f; /* 默认断开状态红点 */
margin-right: 8px;
transition: background-color 0.3s;
}
.status-dot.connected {
background-color: #52c41a; /* 连接成功绿点 */
}
#messageLog {
height: 300px;
border: 1px solid #d9d9d9;
background-color: #fafafa;
padding: 10px;
border-radius: 5px;
overflow-y: auto;
margin-bottom: 15px;
font-family: 'Courier New', Courier, monospace;
font-size: 13px;
}
.log-item {
margin-bottom: 6px;
line-height: 1.4;
}
.log-inbound { color: #1890ff; } /* 收到消息蓝色 */
.log-outbound { color: #52c41a; } /* 发送消息绿色 */
.log-system { color: #8c8c8c; font-style: italic; } /* 系统消息灰色 */
.input-group {
display: flex;
gap: 10px;
}
#inputMsg {
flex: 1;
padding: 10px;
border: 1px solid #d9d9d9;
border-radius: 4px;
outline: none;
font-size: 14px;
}
#inputMsg:focus {
border-color: #40a9ff;
}
button {
padding: 0 20px;
background-color: #1890ff;
color: #fff;
border: none;
border-radius: 4px;
cursor: pointer;
font-size: 14px;
transition: background 0.2s;
}
button:hover {
background-color: #40a9ff;
}
button:disabled {
background-color: #f5f5f5;
color: #rgba(0, 0, 0, 0.25);
border-color: #d9d9d9;
cursor: not-allowed;
}
</style>
</head>
<body>

<div class="container">
<h2>WebSocket 简易客户端</h2>

<div class="status-container">
<div id="statusDot" class="status-dot"></div>
<span id="statusText">正在尝试连接服务器...</span>
</div>

<div id="messageLog"></div>

<div class="input-group">
<input type="text" id="inputMsg" placeholder="请输入要发送给 Netty 的消息..." disabled autocomplete="off">
<button id="sendBtn" onclick="sendMessage()" disabled>发送</button>
</div>
</div>

<script>
// 对应服务端的端口和绑定的路由路径 "/hello"
const serverUrl = "ws://127.0.0.1:8888/hello";
let socket;

const statusDot = document.getElementById("statusDot");
const statusText = document.getElementById("statusText");
const messageLog = document.getElementById("messageLog");
const inputMsg = document.getElementById("inputMsg");
const sendBtn = document.getElementById("sendBtn");

// 初始化连接
function initWebSocket() {
appendLog("正在连接到: " + serverUrl, "system");

// 创建浏览器原生 WebSocket 实例
socket = new WebSocket(serverUrl);

// 1. 监听连接成功事件
socket.onopen = function (event) {
statusDot.classList.add("connected");
statusText.innerText = "已成功连接至 Netty 服务端";
inputMsg.disabled = false;
sendBtn.disabled = false;
appendLog("连接建立成功!可以开始发送消息了。", "system");
};

// 2. 监听收到消息事件
socket.onmessage = function (event) {
appendLog(event.data, "inbound");
};

// 3. 监听连接关闭事件
socket.onclose = function (event) {
statusDot.classList.remove("connected");
statusText.innerText = "与服务端断开连接";
inputMsg.disabled = true;
sendBtn.disabled = true;
appendLog("连接已断开。若要重新连接,请刷新页面。", "system");
};

// 4. 监听发生错误事件
socket.onerror = function (error) {
appendLog("WebSocket 发生错误,请检查服务端是否正常启动。", "system");
};
}

// 发送消息函数
function sendMessage() {
const msg = inputMsg.value.trim();
if (!msg) return;

// 调用原生 API 发送文本内容(Netty 服务端会在 TextWebSocketFrame 中接收到这段字符串)
socket.send(msg);

// 打印本地发送日志并清空输入框
appendLog("【我发送的】" + msg, "outbound");
inputMsg.value = "";
inputMsg.focus();
}

// 辅助函数:向控制面板追加日志行并自动滚动触底
function appendLog(text, type) {
const item = document.createElement("div");
item.className = "log-item log-" + type;

// 添加当前时间戳
const now = new Date().toTimeString().split(' ')[0];
item.innerText = `[${now}] ${text}`;

messageLog.appendChild(item);
messageLog.scrollTop = messageLog.scrollHeight; // 滚动条自动保持在最底部
}

// 允许在输入框中按回车键直接发送
inputMsg.addEventListener("keydown", function(event) {
if (event.key === "Enter") {
sendMessage();
}
});

// 页面加载完毕后自动发起连接
window.onload = initWebSocket;
</script>

</body>
</html>

在 WebSocket 的建立过程中,101 状态码(101 Switching Protocols,切换协议)是整个长连接能够顺利诞生的“入场券”。 如果没有 101 状态码,WebSocket 根本无法启动。简单来说,客户端通过 HTTP 发起升级申请,服务端通过回传 101 状态码正式批准 “变身” 为 WebSocket 协议。为什么要这样设计呢?这是因为:

WebSocket 的设计目标是建立全双工、低延迟的长连接通道,但它在诞生之初面临一个现实问题:如何优雅地穿透互联网上各种复杂的防火墙和网关。如果直接在非标准端口上启动一套全新的 TCP 协议,很容易被各大路由节点当成恶意流量直接墙掉。为了兼容现有的 Web 基础设施,WebSocket 巧妙地选择了借鸡生蛋:

  • 先通过标准的 HTTP 协议(80/443 端口)发起普通的握手请求。
  • 握手一旦达成,立刻在原地将协议“升级切换”为 WebSocket 协议。
  • HTTP 规范中,专门定义了一个状态码来应对这种“中途换道”的场景,这就是 101 Switching Protocols。

我们可以通过客户端与 Netty 服务端在握手瞬间的实际抓包报文,来看 101 是如何起作用的:

  • 阶段一:客户端发起升级申请(HTTP GET)。当你在浏览器执行 new WebSocket(“ws://…”) 时,前端会发出一个特殊的 HTTP 请求,里面包含了“我要升级”的强烈信号:

    1
    2
    3
    4
    5
    6
    7
    GET /hello HTTP/1.1
    Host: 127.0.0.1:8888
    Upgrade: websocket <-- 【信号1】告诉服务器:我想把协议升级为 websocket
    Connection: Upgrade <-- 【信号2】告诉服务器:不要断开,我们要升级连接
    Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
    Origin: null
    Sec-WebSocket-Version: 13
  • 阶段二:服务端同意升级,回传 101(HTTP 响应)。Netty 服务端底层的 WebSocketServerProtocolHandler 接收到这个请求后,进行跨域和路径校验。如果通过,它会立刻原路返回一个 101 状态码的响应:

    1
    2
    3
    4
    HTTP/1.1 101 Switching Protocols  <-- 核心:老板同意了!协议现在正式切换!
    Upgrade: websocket <-- 确认升级为 websocket
    Connection: Upgrade <-- 确认连接保持并升级
    Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=

一旦这个包含 101 的响应成功送达浏览器,原有的 HTTP 协议生命周期在这一瞬间彻底终结。在 Netty 层面,底层的流水线(Pipeline)会发生非常奇妙的变化:

  • 老组件退场:专门处理 HTTP 解析的处理器(如 HttpServerCodec、HttpObjectAggregator)会自动功成身退。因为后续进来的不再是普通的 HTTP 报文,它们留在链条上只会报错。
  • 新组件接管:WebSocketServerProtocolHandler 会把负责 WebSocket 帧解析(Frame Decoder/Encoder)的核心组件动态挂载到管道的最前线。
  • 通道蜕变:此后,这个 TCP 通道上的内容不再有 GET、POST 或 Content-Type,取而代之的是由 Opcode(操作码)和 Payload(有效载荷)组成的 WebSocket 数据帧(Frame)。


加入心跳检测机制

下面我们对之前的 MyWebSocketServer 和业务 Handler 进行升级。规定:如果客户端连续 10 秒没有任何动静(读空闲),服务端就认为其掉线,主动关闭连接。


修改服务器引导类

我们需要在管道中挂载 IdleStateHandler。注意:它必须放在你的自定义业务处理器之前。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// ... 前置的 HttpCodec, ChunkedWriteHandler, HttpObjectAggregator 保持不变 ...

/*
* 核心组件:添加空闲状态检测处理器
* 参数含义:
* 1. readerIdleTimeSeconds: 10 秒内如果服务端没有【读】到客户端的任何数据,触发 READER_IDLE 事件
* 2. writerIdleTimeSeconds: 0 秒(设置为0代表不关心、不检测写空闲)
* 3. allIdleTimeSeconds: 0 秒(设置为0代表不关心、不检测读写空闲)
*/
pipeline.addLast(new IdleStateHandler(10, 0, 0, TimeUnit.SECONDS));

// WebSocket 协议处理器
pipeline.addLast(new WebSocketServerProtocolHandler("/hello"));

// 自定义业务处理器(我们在里面拦截事件)
pipeline.addLast(new MyWebSocketHandler());


修改业务处理器

在 MyWebSocketHandler 增加重写实现方法 userEventTriggered,并适当修改消息响应:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
/**
* 在这里捕捉 Netty 抛出的用户自定义事件(如空闲状态事件)
*/
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception {
if (evt instanceof IdleStateEvent) {
IdleStateEvent event = (IdleStateEvent) evt;

// 判断具体是哪种空闲
if (event.state() == IdleState.READER_IDLE) {
System.out.println("【心跳检测】客户端 ["+ ctx.channel().remoteAddress() +"] 已经连续 10 秒没有任何读写行为,判定为僵尸连接!");

// 通常先向客户端发送一个最后通牒(如特殊的心跳关闭帧),然后主动 close 掉线通道
ctx.channel().close();
}
} else {
// 如果是其他事件,继续向后传递
super.userEventTriggered(ctx, evt);
}
}

@Override
protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception {
// 1. 获取客户端传过来的文本内容
String text = msg.text();
System.out.println("【WebSocket服务端】收到来自客户端[" + ctx.channel().remoteAddress() + "]的消息: " + text);

// 2. 回复客户端一条消息
String responseText = "PING".equals(text) ? "PONG" : "【服务器回复】" + LocalDateTime.now() + " -> 我已收到你的消息: " + text;
ctx.channel().writeAndFlush(new TextWebSocketFrame(responseText));
}


升级 HTML 客户端 - 建立保活机制

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
let socket;
let heartbeatIntervalId; // 用来存放定时器 ID

function initWebSocket() {
socket = new WebSocket("ws://127.0.0.1:8888/hello");

socket.onopen = function (event) {
// ... 原有的 UI 解除置灰逻辑 ...
appendLog("连接建立成功!", "system");

// 👉 开启客户端心跳保活定时器:每隔 5 秒向 Netty 发送一个特制的心跳字符串
startHeartbeat();
};

socket.onmessage = function (event) {
// 过滤掉服务端返回的心跳回应,不让它污染用户的聊天日志
if (event.data === "PONG" || event.data.includes("心跳")) {
return;
}
appendLog(event.data, "inbound");
};

socket.onclose = function (event) {
// ... 原有的置灰逻辑 ...
appendLog("连接已断开。", "system");

// 👉 连接断开了,必须清除定时器,防止死循环空发
stopHeartbeat();
};
}

// 👉 启动心跳
function startHeartbeat() {
heartbeatIntervalId = setInterval(() => {
if (socket && socket.readyState === WebSocket.OPEN) {
// 发送两端约定好的心跳帧文本,通常叫 "PING"
socket.send("PING");
// 本地可以打个淡灰色的日志,证明客户端还在努力活着
console.log("[心跳] 向服务端发送 PING 保活...");
}
}, 5000); // 5秒发一次,远小于服务端的 10 秒限制,安全
}

// 👉 停止心跳
function stopHeartbeat() {
if (heartbeatIntervalId) {
clearInterval(heartbeatIntervalId);
}
}


支持连接鉴权

在实际生产环境中,连接鉴权(Authentication)是长连接系统的第一道安全防线。如果允许匿名用户随意建立 WebSocket 连接,服务器的内存和网络带宽很容易被恶意攻击者榨干。要实现像 ws://127.0.0.1:8888/hello?token=abcdefg&userId=1001 这样带参数的握手并在握手阶段进行安全校验,我们需要利用 Netty 责任链的特殊生命周期。


核心原理:利用 HTTP 阶段拦截鉴权

前面我们说到,WebSocket 在正式建立前,必须经历一个 HTTP 101 升级握手 阶段。这意味着:带参数的 URL 请求,本质上是一个标准的 HTTP GET 请求。

在 Netty 中,WebSocketServerProtocolHandler 一旦检测到匹配的路径(如 /hello),就会自动处理 101 升级。因此,我们必须在升级发生之前,拦截这个 HTTP 请求,解析出里面的 token 和 userId 并进行合法性校验。 如果校验失败,直接返回 HTTP 401 Unauthorized 并斩断连接,根本不给它升级成 WebSocket 的机会。


具体的实现

编写鉴权处理器

我们编写一个专门的 HTTP 入站处理器 AuthHandler。它继承自 ChannelInboundHandlerAdapter(处理普通的 HTTP 对象)。

在实际的分布式大型项目中,握手阶段的 token 通常建议使用 JWT(JSON Web Token)。JWT 中本身就加密包裹了 userId、过期时间以及权限标识。AuthHandler 拦截到请求后,直接在内存中用秘钥解密 JWT。这样不需要查数据库或 Redis,就能在 0.1 毫秒内完成鉴权,极大减轻了长连接握手风暴(如服务器重启时百万客户端同时重连)对后台核心数据库的冲击!

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
import io.netty.channel.ChannelFutureListener;
import io.netty.channel.ChannelHandlerContext;
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.handler.codec.http.*;
import io.netty.util.AttributeKey;
import lombok.extern.slf4j.Slf4j;
import java.util.List;
import java.util.Map;

@Slf4j
public class AuthHandler extends ChannelInboundHandlerAdapter {

// 定义一个 Channel 的属性 Key,用于在同一次连接的后续 Handler 中共享用户 ID
public static final AttributeKey<String> USER_ID_KEY = AttributeKey.valueOf("userId");

@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
// 1. 检查是否是 HTTP 全量请求对象(握手阶段的数据)
if (msg instanceof FullHttpRequest) {
FullHttpRequest request = (FullHttpRequest) msg;
String uri = request.uri();

// 2. 检查是否是发往我们 WebSocket 路由的请求
if (uri.startsWith("/hello")) {
// 3. 解析 URL 中的参数(使用 Netty 自带的 QueryStringDecoder)
QueryStringDecoder decoder = new QueryStringDecoder(uri);
Map<String, List<String>> parameters = decoder.parameters();

// 提取参数
List<String> tokenList = parameters.get("token");
List<String> userIdList = parameters.get("userId");

String token = (tokenList != null && !tokenList.isEmpty()) ? tokenList.get(0) : null;
String userId = (userIdList != null && !userIdList.isEmpty()) ? userIdList.get(0) : null;
log.info("【鉴权层】收到握手申请,userId: {}, token: {}", userId, token);

// 4. 执行安全校验
if (!validateToken(userId, token)) {
log.warn("【鉴权层】拒绝连接!用户 [{}] 鉴权失败。", userId);
// 鉴权失败,返回 HTTP 401,拒绝升级
sendUnauthorizedAndClose(ctx);
return; // 斩断责任链,不再向下传递
}

// 5. 鉴权成功:将用户 ID 绑定到当前 Channel 的属性中,后续聊天时就能直接知道“我是谁”
ctx.channel().attr(USER_ID_KEY).set(userId);

/*
* 6. 因为 URL 后面带了参数(如 /hello?token=...),而底层的
* WebSocketServerProtocolHandler 默认是【精确匹配】字符串 "/hello" 的。
* 如果不把参数裁剪掉,后面的 WebSocket 升级组件会报 404!
*/
request.setUri(decoder.path());
}
}

// 鉴权通过,继续向后传递(交给后面的 HttpAggregation 和 WebSocket 协议处理器)
super.channelRead(ctx, msg);
}

/**
* 模拟安全鉴权逻辑(实际生产中这里对接 Redis、JWT 校验或数据库)
*/
private boolean validateToken(String userId, String token) {
if (userId == null || token == null) {
return false;
}
// 举例:假设正确的 token 就是 "abcdefg"
return "abcdefg".equals(token);
}

/**
* 拒绝策略:返回 401 状态码并关闭连接
*/
private void sendUnauthorizedAndClose(ChannelHandlerContext ctx) {
FullHttpResponse response = new DefaultFullHttpResponse(
HttpVersion.HTTP_1_1,
HttpResponseStatus.UNAUTHORIZED
);
response.headers().set(HttpHeaderNames.CONTENT_TYPE, "text/plain; charset=UTF-8");
HttpUtil.setContentLength(response, 0);

// 异步发送响应,并在发送完毕后立刻强行关闭 TCP 连线
ctx.writeAndFlush(response).addListener(ChannelFutureListener.CLOSE);
}
}


修改服务端流水线

在 MyWebSocketServer 中,我们将 AuthHandler 塞进 Pipeline。它的位置必须放在 HTTP 编解码和聚合器之后,但在 WebSocketServerProtocolHandler 之前!

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
@Override
protected void initChannel(SocketChannel ch) throws Exception {
ChannelPipeline pipeline = ch.pipeline();

pipeline.addLast(new HttpServerCodec());
pipeline.addLast(new ChunkedWriteHandler());
pipeline.addLast(new HttpObjectAggregator(8192));
pipeline.addLast(new IdleStateHandler(10, 0, 0, TimeUnit.SECONDS));

// 👉🏻 核心:在正式升级 WebSocket 之前,切入我们的鉴权处理器
pipeline.addLast(new AuthHandler());

pipeline.addLast(new WebSocketServerProtocolHandler("/hello"));
pipeline.addLast(new MyWebSocketHandler());
}


在业务层获取用户

当连接成功升级后,我们在聊天业务 MyWebSocketHandler 里就不需要每次都带 token 了。可以直接从当前 Channel 里取出刚才在 AuthHandler 里绑定的 userId:

1
2
3
4
5
6
7
8
9
10
11
@Override
protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception {
// 1. 获取客户端传过来的文本内容
String userId = ctx.channel().attr(AuthHandler.USER_ID_KEY).get();
String text = msg.text();
System.out.println("【WebSocket服务端】收到用户[" + userId + "]的消息: " + text);

// 2. 回复客户端一条消息
String responseText = "PING".equals(text) ? "PONG" : "【服务器回复】" + LocalDateTime.now() + " -> [用户 " + userId + " ] 我已收到你的消息: " + text;
ctx.channel().writeAndFlush(new TextWebSocketFrame(responseText));
}


升级 HTML 客户端

修改 HTML 里的 serverUrl,带上参数去连接。你可以分别尝试正确参数和错误参数:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
// 实验 1:使用正确合法的 token
const serverUrl = "ws://127.0.0.1:8888/hello?token=abcdefg&userId=1001";

// 实验 2:修改 token 模拟黑客攻击(将其改为 "wrong_token")
// const serverUrl = "ws://127.0.0.1:8888/hello?token=wrong_token&userId=1001";


// 监听连接关闭事件
socket.onclose = function (event) {
statusDot.classList.remove("connected");

// 默认的系统提示
let reasonText = "与服务端断开连接";

/**
* 利用 event.code 状态码来推断具体的错误原因!
* 虽然 error 里没写,但关闭码里写得清清楚楚
*/
if (event.code === 1006) {
// 1006 代表异常关闭(非正常四次挥手断开)
reasonText = "连接失败:服务端拒绝了握手请求(可能未启动、路径错误、或鉴权失败)";
appendLog(`[错误代码: ${event.code}] 握手被拒绝,请检查 token 是否正确、服务端是否存活。`, "system");
} else {
appendLog(`[连接关闭] 关闭代码: ${event.code}, 原因: ${event.reason || '未指定'}`, "system");
}

statusText.innerText = reasonText;
inputMsg.disabled = true;
sendBtn.disabled = true;
};

// 监听发生错误事件
socket.onerror = function (error) {
// 直接打印 error 拿不到具体字符串,但通过控制台对象展开可以看到其类型
console.error("【WebSocket 底层错误触发】:", error);
appendLog("【网络提示】浏览器底层检测到 WebSocket 异常,连接即将关闭...", "system");
};

实验1:

实验2:


群聊与私聊的实现

在原生的 NIO 中,为了管所有客户端的 SocketChannel,我们通常得自己写一个 ConcurrentHashMap 来维护用户 ID 和 Channel 的映射关系。到了 Netty 中,框架为我们提供了一个极其强大的容器 —— ChannelGroup。它可以帮我们把成千上万的 Channel 组织成一个聚合组。不仅能自动管理 Channel 的生命周期(当客户端断开连接时,Netty 会自动将其从 ChannelGroup 中移除,不需要你手动 remove),还提供了优雅的、批量的高性能异步广播机制。下面我们结合前面的 AuthHandler(已将 userId 绑定到 Channel 属性中),来实现一个支持精准单聊(私聊)和全员广播(群聊)的案例。


编写全网连接管理器

在实际开发中,我们通常会抽象出一个单例的 SessionManager,用来同时兼顾 “全员群聊(ChannelGroup)” 和 “精准私聊(ConcurrentHashMap)”。

MySessionManager:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
import io.netty.channel.Channel;
import io.netty.channel.group.ChannelGroup;
import io.netty.channel.group.DefaultChannelGroup;
import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;
import io.netty.util.concurrent.GlobalEventExecutor;
import lombok.extern.slf4j.Slf4j;
import java.util.concurrent.ConcurrentHashMap;

@Slf4j
public class MySessionManager {

// 1. 全局单例
private MySessionManager() {}
private static final MySessionManager INSTANCE = new MySessionManager();
public static MySessionManager getInstance() { return INSTANCE; }

// 2. ChannelGroup:用于管理所有在线的连接,专门负责【全员群聊广播】。GlobalEventExecutor.INSTANCE 是一个全局单线程的事件执行器,用来驱动群组内的异步操作
private final ChannelGroup globalGroup = new DefaultChannelGroup(GlobalEventExecutor.INSTANCE);

// 3. Map:用于维护 用户ID -> Channel 的映射,专门负责【特定用户精准私聊】
private final ConcurrentHashMap<String, Channel> userChannelMap = new ConcurrentHashMap<>();


/**
* 当用户鉴权通过且连接建立成功时,注册会话
*/
public void register(String userId, Channel channel) {
// 自动加入群组(如果已存在会自动忽略,内部是 Set 结构)
globalGroup.add(channel);
// 绑定用户私聊映射
userChannelMap.put(userId, channel);
log.info("【会话管理】用户 [{}] 登录上线,当前在线人数: {}", userId, userChannelMap.size());
}

/**
* 当用户断开连接或心跳超时被踢时,注销会话
*/
public void unregister(Channel channel) {
// 从 Channel 属性中取出绑定的 userId
String userId = channel.attr(AuthHandler.USER_ID_KEY).get();
if (userId != null) {
userChannelMap.remove(userId);
}
// 提示:其实 globalGroup.remove(channel) 这里不写也行,Netty 发现 Channel close 后会自动从组里踢掉它
globalGroup.remove(channel);
log.info("【会话管理】用户 [{}] 离线,当前在线人数: {}", userId, userChannelMap.size());
}

/**
* 功能一:全员广播(群聊)
*/
public void broadcast(String message) {
// ChannelGroup 的神奇之处:直接把对象丢进去,它会异步向组内所有的 Channel 发送,效率极高!
// 注意:WebSocket 必须广播 TextWebSocketFrame 对象
globalGroup.writeAndFlush(new TextWebSocketFrame("【系统广播】" + message));
}

/**
* 功能二:精准私聊(单聊)
* @param toUserId 接收方用户 ID
* @param message 消息内容
*/
public void sendToUser(String toUserId, String message) {
Channel targetChannel = userChannelMap.get(toUserId);
if (targetChannel != null && targetChannel.isActive()) {
// 找到通道,精准推送
targetChannel.writeAndFlush(new TextWebSocketFrame(message));
} else {
log.warn("【会话管理】发送失败:用户 [{}] 当前不在线或连接已失效。", toUserId);
}
}
}


业务处理器升级 - 约定聊天协议以及用户上线离线的绑定

为了让服务端知道这条消息是“群发”还是“私聊给某人”,我们规定一个最简单的文本协议格式(通常生产环境用 JSON,这里用文本切分演示):

  • 群聊格式:ALL:大家好呀!
  • 私聊格式:TO:1002:嘿,哥们,你在吗?

修改:MyWebSocketHandler

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
public class MyWebSocketHandler extends SimpleChannelInboundHandler<TextWebSocketFrame> {

/**
* 在这里捕捉 Netty 抛出的用户自定义事件
*/
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception {
if (evt instanceof IdleStateEvent) {
IdleStateEvent event = (IdleStateEvent) evt;
if (event.state() == IdleState.READER_IDLE) {
System.out.println("【心跳检测】客户端 [" + ctx.channel().remoteAddress() + "] 已经连续 10 秒没有任何读写行为,判定为僵尸连接!");
ctx.channel().close();
}
} else if (evt instanceof WebSocketServerProtocolHandler.HandshakeComplete) { // 👉🏻 处理 WebSocket 握手成功事件(此时协议已升级,且 AuthHandler 注入的 userId 必然存在!)
String userId = ctx.channel().attr(AuthHandler.USER_ID_KEY).get();
if (userId != null) {
// 此时注册,安全无误!
MySessionManager.getInstance().register(userId, ctx.channel());
// 全员广播上线提示
MySessionManager.getInstance().broadcast("欢迎用户 [" + userId + "] 进入聊天室!");
}
} else {
// 如果是其他事件,继续向后传递
super.userEventTriggered(ctx, evt);
}
}

/**
* 为了让服务端知道这条消息是“群发”还是“私聊给某人”,我们规定一个最简单的文本协议格式(通常生产环境用 JSON,这里用文本切分演示):
* 群聊格式:ALL:大家好呀!
* 私聊格式:TO:1002:嘿,哥们,你在吗?
*/
@Override
protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception {
String currentUserId = ctx.channel().attr(AuthHandler.USER_ID_KEY).get();
String rawText = msg.text();

// 过滤心跳
if ("PING".equals(rawText)) {
ctx.channel().writeAndFlush(new TextWebSocketFrame("PONG"));
return;
}
System.out.println("【业务层】收到用户 ["+ currentUserId +"] 的原始消息: " + rawText);

// 解析约定的协议
if (rawText.startsWith("ALL:")) {
// 群聊逻辑
String content = rawText.substring(4);
String broadcastMsg = String.format("【群聊】用户[%s]说: %s", currentUserId, content);
MySessionManager.getInstance().broadcast(broadcastMsg);

} else if (rawText.startsWith("TO:")) {
// 私聊逻辑,格式如:TO:1002:内容
String[] parts = rawText.split(":", 3);
if (parts.length == 3) {
String toUserId = parts[1];
String content = parts[2];

String privateMsg = String.format("【私聊】用户[%s]悄悄对你说: %s", currentUserId, content);
MySessionManager.getInstance().sendToUser(toUserId, privateMsg);

// 给发送者自己也回显一条,体验更好
ctx.channel().writeAndFlush(new TextWebSocketFrame(String.format("【你对%s说】: %s", toUserId, content)));
} else {
ctx.channel().writeAndFlush(new TextWebSocketFrame("【系统提示】格式错误,私聊请使用: TO:用户ID:内容"));
}
} else {
ctx.channel().writeAndFlush(new TextWebSocketFrame("【系统提示】不认识的指令。群发请输入 ALL:内容,私聊请输入 TO:用户ID:内容"));
}
}

/**
* 砍掉原有的 channelActive 注册逻辑!
* 因为 Active 时 HTTP 数据还没读进来,userId 必然为 null。
*/
@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
// 保持现状或只打印底层连接,不在这里写业务会话绑定
System.out.println("【系统提示】底层通道已激活: " + ctx.channel().remoteAddress());
}

@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception {
String userId = ctx.channel().attr(AuthHandler.USER_ID_KEY).get();
// 离线注销
MySessionManager.getInstance().unregister(ctx.channel());
if (userId != null) {
MySessionManager.getInstance().broadcast("用户 [" + userId + "] 离开了聊天室。");
}
}

@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
System.err.println("【WebSocket服务端】发生异常" + cause);
ctx.close(); // 关闭通道
}
}

注意,上述 MySessionManager.getInstance().register 必须在 userEventTriggered 才可以(而不是 channelActive),这是网络 TCP 连接的建立过程以及Netty 底层的事件流向决定的。简单来说:

  • channelActive 代表的是 “物理通道(TCP)通了”;
  • 而 userEventTriggered 里的 HandshakeComplete 代表的是 “业务协议(WebSocket)谈成了”。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
[步骤 1] 客户端发起 TCP 三次握手 

[步骤 2] 🤝 TCP 握手成功!操作系统的 Socket 通道建立

└─► 🔥 触发 channelActive()
(此时:通道里空空如也,客户端连 1 个字节的 HTTP 握手数据都还没发过来)

[步骤 3] 📨 浏览器发送 HTTP 请求报文(GET /hello?userId=1001&token=...)

[步骤 4] ⚙️ 数据流入 Pipeline:HttpServerCodec ──► AuthHandler (解析出 userId 并解析绑定)

[步骤 5] 🎯 WebSocketServerProtocolHandler 进行 101 状态码协议升级响应

└─► 🔥 触发 userEventTriggered() 并且携带 HandshakeComplete 事件
(此时:WebSocket 长连接正式宣告诞生)


联调测试

现在你可以同时打开两个不同的浏览器窗口(或者一个是正常模式,一个是无痕模式),模拟两个不同的用户同时在线:

  • 玩家 A 上线:

    • 修改 HTML 的 URL 里的 userId 为 1001:
      const serverUrl = “ws://127.0.0.1:8888/hello?token=abcdefg&userId=1001”;
    • 打开网页 A。系统日志会显示:[系统广播] 欢迎用户 [1001] 进入聊天室!
  • 玩家 B 上线:

    • 修改同一个 HTML 里的 userId 为 1002:
      const serverUrl = “ws://127.0.0.1:8888/hello?token=abcdefg&userId=1002”;
    • 打开网页 B。此时网页 A 和网页 B 都会同时收到:[系统广播] 欢迎用户 [1002] 进入聊天室!(证明 ChannelGroup 广播成功!)
  • 测试群聊:

    • 在用户 1001 的输入框输入:ALL:大家好,我是1001!
    • 结果:两个网页同时收到了这行群发。
  • 测试精准私聊:
    • 在用户 1001 的输入框输入:TO:1002:哥们,单独聊聊,今晚去吃烧烤吗?
    • 观察结果:
      • 用户 1002 的网页上:弹出了私聊字样:【私聊】用户[1001]悄悄对你说: 今晚去吃烧烤吗?
      • 其他任何不相干用户的网页上:一片安静,什么都收不到。这就是精准定向推送。


用户上线:

广播消息:

单独私聊:


性能深度优化

在高并发的 WebSocket 推送系统中(如百万级在线的聊天室、直播弹幕、即时股票行情),内存暴涨和 OOM(内存溢出) 往往是除了线程阻塞之外最致命的杀手。Netty 为了榨干网络 IO 的性能,默认使用了池化的堆外内存(PooledByteBufAllocator)。堆外内存不受 JVM 垃圾回收(GC)的管辖,如果开发者在使用、跨线程传递过程中少释放了一次,就会造成底层的内核内存泄漏;如果多释放了一次,则会引发严重的 IllegalReferenceCountException: refCnt: 0 崩溃。针对以上高并发推送场景,我们将从 TextWebSocketFrame 的复用跨线程派发 两个维度,深度剖析如何避开内存泄漏的坑。


跨线程池派发任务时的引用计数坑

在 Netty 中,所有的网络出入站对象(如 ByteBuf、TextWebSocketFrame 等)都实现了 ReferenceCounted(引用计数)接口。其核心铁律是:

  • 对象的初始引用计数为 1。
  • 只要引用计数变为 0,对象占用的底层堆外内存就会被立刻释放并回收到对象池。
  • 谁消费,谁释放。 Netty 的 SimpleChannelInboundHandler 会在 channelRead0 结束后自动帮你把 msg 的引用计数减 1。

经典的跨线程泄漏死锁的例子:很多开发者在收到消息后,为了不阻塞 IO 线程(EventLoop),会把消息提交到自定义的业务线程池中处理:

1
2
3
4
5
6
7
8
9
10
11
// ❌ 错误示范:会导致严重的内存泄漏或运行时崩溃
@Override
protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) {
// 将任务提交给业务线程池
businessThreadPool.execute(() -> {
// 这里的 msg 极有可能已经被 IO 线程在 channelRead0 结束时自动 release 掉了!
// 此时业务线程去读取 msg.text() 会直接抛出 refCnt: 0 异常。
String text = msg.text();
System.out.println(text);
});
}

正确的跨线程传递姿势(保留引用):如果必须把包装了 ByteBuf 的对象(如 TextWebSocketFrame)传给其他线程,你必须显式调用 retain() 来增加引用计数,告诉 Netty:“我还要用,你先别回收!”。并在业务线程消费完毕后手动 release()。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// ✅ 正确示范:跨线程传递时手动维护引用计数
@Override
protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) {
// 1. 显式增加引用计数(refCnt 变为 2)
msg.retain();

businessThreadPool.execute(() -> {
try {
String text = msg.text();
// 执行你的核心业务逻辑...
} finally {
// 2. 业务线程消费完毕,必须手动释放(refCnt 减 1,变回 1。随后 IO 线程执行完减 1 变 0 彻底回收)
msg.release();
}
});
}


高并发下群发推送时 TextWebSocketFrame 极致复用

假设一个聊天室里有 1 万名在线用户,有人发了一条消息,服务端需要同时推送给这 1 万个人。

❌ 愚蠢的做法:循环 new 对象。这种做法在高并发下会产生大量的内存碎片,频繁触发池化分配器的锁竞争,甚至导致 OOM。

1
2
3
4
// 循环 1 万次,new 了 1 万个对象,在堆外申请了 1 万次内存,带来巨大的多线程内存争用
for (Channel ch : globalGroup) {
ch.writeAndFlush(new TextWebSocketFrame("大家好!"));
}

✅ 高性能复用方案:使用 ChannelGroup 或手动 retain()。Netty 的 ChannelGroup.writeAndFlush(object) 内部已经做过了极致的复用优化:它不会为每个通道单独复制一份数据,而是通过共享同一个底层 ByteBuf 内存,并对引用计数进行递增的方式实现的。如果你因为特殊的路由逻辑(比如按照用户标签推送),无法使用 ChannelGroup,必须自己写 for 循环群发,那么你必须采用 “单对象 + 手动 retain” 的零拷贝复用策略:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
public void broadcastToUsers(List<Channel> targetChannels, String message) {
if (targetChannels == null || targetChannels.isEmpty()) return;

// 1. 只创建【一个】WebSocket 帧对象(此时 refCnt = 1)
TextWebSocketFrame frame = new TextWebSocketFrame(message);

// 2. 计算需要发送的目标总数
int targetCount = targetChannels.size();

// 3. 关键:除了最后一个通道,前面的 count-1 个通道发送前,都必须 retain() 一次
// 因为每次调用 writeAndFlush,Netty 在底层的尾部处理器(TailContext)中都会自动对该对象调用一次 release()。
// 如果不提前 retain,第一个通道发完,引用计数归零,后面所有的通道发送时都会报错!
for (int i = 0; i < targetCount - 1; i++) {
frame.retain(); // 增加引用计数
}

// 4. 开始群发(所有人共享同一个底层的堆外内存 ByteBuf)
for (Channel ch : targetChannels) {
if (ch.isActive()) {
ch.writeAndFlush(frame);
} else {
// 万一中途发现这个通道已经失效了,Netty 不会帮我们对这个失效通道做 write 之后的自动 release
// 我们必须手动补齐,把刚才预先加的引用扣掉,防止内存泄漏!
frame.release();
}
}
}


监控和排查内存泄漏

即使代码写得再小心,高并发下依然可能存在链路遗漏导致的泄漏。Netty 提供了工业级的内存泄漏采样探测器。在生产环境或本地压力测试启动时,通过 JVM 参数开启最高级别的泄漏检测:

1
2
3
4
5
# DISABLED: 关闭检测(不推荐)。
# SIMPLE: 简易级别,占用 1% 采样率。如果发生泄漏,会打印一行警告日志,但没有堆栈信息。
# ADVANCED: 高级级别,高采样率,发生泄漏时会记录该 ByteBuf 是在哪个类的哪个方法被创建或被操作的(非常实用!)。
# PARANOID: 偏执级别,100% 全量采样(压测和排查期必开,生产环境视 CPU 负载而定)。
-Dio.netty.leakDetection.level=PARANOID

如果发生泄漏,控制台会弹出类似如下的恐怖警告:

1
2
3
4
5
LEAK: ByteBuf.release() was not called before it's garbage-collected. 
See https://netty.io/wiki/reference-counted-objects.html for more information.
Recent access records: 4
#1: com.owlias.server.zdemo.AuthHandler.channelRead(AuthHandler.java:42)
#2: io.netty.handler.codec.MessageToMessageDecoder.channelRead(...)

一旦看到这个日志,不要慌,顺着它提示的 #1、#2 堆栈,去找最后一次创建或调用 retain(),但后续没有安全闭环 release() 的代码位置,将其放进 try-finally 块中确保释放即可。